
# 1、创建flink执行环境
env = StreamExecutionEnvironment.get_execution_environment()

# 修改并行度
env.set_parallelism(1)

# 2、读取数据,得到DataStream，相当于RDD（有界流）
scores_ds = env.read_text_file("G:\LanZhiPeiXun\Flink\data\score.txt")

# 方法1、使用lambda表达式的方式，每一条数据都需要创建数据库连接，吞吐量会很低
def map_fun(score):
    id = score.split(",")[0]

    # 1、创建数据库连接
    con =pymysql.connect(host="master",part=3306)










env.execute()